昨天 D11 介紹完晚物化,WHERE g = 5 拿到一張「這幾列符合條件」的位圖,還沒真的把 payload 解出去。那輸出前還要跟 delete 交集嗎?這就是今天要挖的這一截~
今天換個講法,這次從最底層開始,先不管那些名詞
在 MySQL 這種資料庫裡,資料存在資料庫自己管的檔案裡,你下 DELETE,它就自己去把那一列改掉,你不用管細節。
資料湖不是這樣。資料湖的一張表,說穿了就是雲端硬碟上的一堆檔案。
orders 這張表實際上在 S3 上長這樣:
s3://mybucket/orders/
part-00001.parquet ← 裡面裝 100 萬列訂單
part-00002.parquet ← 再 100 萬列
part-00042.parquet ← 再 100 萬列
...
沒有資料庫伺服器在跑,沒有背景程式在守,就是一堆躺在雲端的檔案。
Spark 想查資料就是自己把這些檔案下載下來讀。
這件事先記住,因為接下來所有的麻煩都是從這裡長出來的。
客人打電話來說訂單 999 要取消,你下:
DELETE FROM orders WHERE id = 999;
MySQL 那邊,資料庫程式打開檔案、找到那一行、蓋掉、關檔,一氣呵成。
資料湖沒有這個「資料庫程式」,只有 Spark 這種讀寫端跟 S3 上一堆 Parquet 檔。而 Parquet 一個致命限制,寫出去之後不能改。
「不能改」不是 Parquet 龜毛,是它為了「一次寫、多次讀」做的取捨。Column chunk、page、footer 都是壓好、算好、鎖好位置的,中間插一個「這列刪了」進去,後面所有 offset 都要重算,等於把整個檔案重寫。而且物件儲存本身也沒有「就地改一個 byte」這種操作,只能 PUT 整份檔案覆蓋。
於是問題變成,「這個檔案的第 12345 列不算數」這件事,既然沒辦法寫回檔案裡,那要寫到哪裡去?
整篇文章就是在回答這一句。
Copy-on-write 的做法很笨但很直白:
part-00042.parquet 整份下載(100 萬列)part-00099.parquet(999,999 列)好處是讀的人完全不用做任何額外的事,打開檔案讀到的就是正確答案,跟以前一樣。
壞處是為了刪 1 列,重寫了 128 MB。
一天刪一次還可以忍,可是資料是從 Kafka 一直流進來、每分鐘就有一筆取消訂單的話,你每分鐘重寫一次 128 MB,一天重寫 1440 次,硬碟跟網路直接被打死。
所以 COW 只在「很少改」的表上划算,例如純批次寫的 fact table。
CDC upsert 這種情境要換另一條路。
Merge-on-read 的想法很生活,檔案不動,另外拿一張便條記「這個檔案的第 X 列不算數」,讀的人邊讀檔案邊對便條,看到那列就自動跳過。
「便條」在 Iceberg 裡就叫 delete file,本身也是一個 Parquet 檔,跟資料檔一起掛在 manifest 上,manifest 那一列會標「這是 delete file,不是 data file」。這是 Iceberg V2 的主場,也是今天的主題。
這是 V3 才有的路,把 MOR 的便條從「一列一列寫」換成「一張跟資料檔一樣長的位圖」。放到第六步再細講,因為它是路線二的改良版,要先懂路線二的痛才知道它在改什麼。
MOR 說要寫便條,可是便條上到底寫什麼?有兩種寫法。
便條內容:
(檔案 = part-00042.parquet, 第 12345 列)
讀的人很好處理。 因為讀檔案本來就是從第 0 列、第 1 列、第 2 列一路往下讀,便條上的位置也事先排好序,所以兩邊各放一根手指一起往下走:
資料讀到第 ... 12343 12344 12345 12346 ...
便條上寫著 ↓
12345
對上了,這列跳過
手指只往前不往後,走一遍就結束,成本幾乎是零。這個動作有個名字,叫 sorted merge(兩份都排好序的東西一起走一遍合起來)。
但寫的人有麻煩,他要先知道「第 12345 列」這個數字。
要下的 SQL 是 WHERE id = 999,這句話裡沒有「位置」這個東西,只有一個主鍵值。所以寫入端要先跑一輪 lookup,把 id = 999 對應到「part-00042 的第 12345 列」,才寫得下這張便條。
Batch job 不太痛,反正它本來就會全掃一次;Streaming 就痛,來源那邊只丟主鍵,還要繞回去查一次位置。
便條內容:
(欄位 = id, 值 = 999)
翻成人話,「所有 id = 999 的都不算數,不管它們住在哪個檔案的第幾列」。
寫的人很好處理。 來源給什麼就寫什麼,來源說 id = 999 那便條就寫 id = 999,不用先跑 lookup。這對 CDC 是天生的形狀。
但讀的人麻煩。 因為便條沒說位置,讀端要自己找出「哪幾列的 id 是 999」:
{999, 1024, 3311, ...}
100 萬列就是 100 萬次查詢。這個動作叫 anti-join(反向比對,比對上的要丟掉),比寫法 A 的「手指走一遍」貴一個等級。
而且還有一個副作用,你的查詢明明是 SELECT amount FROM orders WHERE g = 5,根本沒用到 id 這欄,可是為了比對刪除,你還是得把 id 這一整欄解出來,憑空多解一欄。
| 便條寫什麼 | 誰辛苦 | |
|---|---|---|
| Positional delete | 42 號檔第 12345 列 | 寫的人(要先找位置) |
| Equality delete | id = 999 | 讀的人(每列都要查一次) |
同一件苦工(把主鍵對應到位置),只是決定推給誰做。
因為它是「為什麼需要寫法 B」的原因。
CDC = Change Data Capture,變更資料擷取。
情境:公司正式的訂單資料在 MySQL 裡(線上交易系統),可是分析用的資料湖也想同步這份資料好跑報表。你不會希望每次都全表 dump 一次,太慢,所以改成「MySQL 每有一筆增刪改就推一則變更事件到 Kafka,資料湖那邊訂閱 Kafka,把事件套進資料湖」。這條 Kafka 上流動的資料就叫 CDC event。
一則典型的 CDC event 長這樣:
{op: DELETE, table: orders, id: 999}
Kafka 只知道主鍵,不知道也不會知道「id = 999 那筆住在 part-00042 的第 12345 列」,因為它跟 Parquet 本來就沒關係。這時候寫入端如果要寫 positional delete,就得繞回資料湖查一次位置,慢又貴;改寫 equality delete 直接把主鍵抄下來就好。這就是為什麼 equality delete 存在。
現在把整條讀取流程走一遍,假設一個小檔案只有 10 列(真實是幾十萬列,這裡縮小方便看):
列號 g amount
0 3 100
1 5 200
2 5 150
3 1 900
4 5 300
5 5 250 ← 這列 id=999,上禮拜被刪了
6 2 400
7 7 800
8 5 120
9 4 600
查詢:SELECT amount FROM orders WHERE g = 5
Step 1:只解 g 這一欄,跑條件
不要碰 amount,因為還不知道要不要用它。只解 g,算出一張「誰符合」的名單,這張名單就是一排 0 跟 1:
列號: 0 1 2 3 4 5 6 7 8 9
v_pred: 0 1 1 0 1 1 0 0 1 0
這就是 D11 講的晚物化跟 selection vector,v_pred 是「符合 WHERE 的名單」。
Step 2:讀便條,算出「誰被刪了」的名單
去讀綁在這個檔案上的 delete file,算出另一排 0 跟 1:
v_delete: 0 0 0 0 0 1 0 0 0 0
第 5 列因為 delete file 裡寫著「id = 999 被刪」而且它就是 id = 999,標 1,其他都 0。
Step 3:兩張名單合起來
「符合條件」還要「還活著」:
v_pred: 0 1 1 0 1 1 0 0 1 0
v_delete: 0 0 0 0 0 1 0 0 0 0
v_out: 0 1 1 0 1 0 0 0 1 0 ← v_pred AND NOT v_delete
第 5 列本來符合 g = 5(v_pred 是 1),可是被刪了(v_delete 是 1),最後 v_out 變成 0。
這就是 v_out = v_pred AND NOT v_delete 的意思,AND 是「兩邊都要是 1 才是 1」、NOT 是「反過來」。純粹的位元運算,CPU 做這個超快,D2 那條 vectorized loop 天生就是為了這個而生的。
Step 4:拿最後名單去解 amount
只解第 1、2、4、8 列的 amount,第 5 列連解都不用解,直接省下來。
「delete 沒有改變晚物化的形狀」是什麼意思? 晚物化的精神是「先用便宜的方式算出一張名單,最後才根據名單去解貴的資料」。加了刪除之後,名單的算法從一步變兩步(一張條件名單、一張刪除名單、合起來),但**「先算名單,最後才解資料」這個順序沒變**。刪除是插隊在名單定案之前,不是「先把資料全解出來再把刪掉的丟掉」。整個架構沒被破壞,只多了一次 bitmap 運算。
便條的痛點
一份資料檔可能綁一堆 delete file,讀時全部要載入、全部要合起來。
Streaming 情境每分鐘寫一次 delete,一天就是 1440 張便條貼在同一批資料檔上,讀時 merge 的路數爆炸。你查一份表,讀路徑上可能要打開 4000 個檔(1 個 data file + 3999 張便條),S3 GET 一次 100 ms,光開檔就把 latency 吃光。讀的成本來自便條的張數,不是刪除的筆數。
當然可以定期做 compaction 把便條合併起來,可是這變成一場永遠追不完的賽跑,跑太少便條堆積讀就慢,跑太多 compaction 自己要重寫檔案,跟寫入端搶頻寬,你被綁在一根鋼索上。
Deletion vector 的解法
改變表達方式。不再是「每次刪除寫一張便條」,而是「一個資料檔永遠只配一張位圖」。
100 萬列的檔案就對應一張 100 萬個 bit 的位圖,第 12345 個 bit 是 1 代表第 12345 列被刪了。
今天刪 3 列,把 3 個 bit 設成 1;明天刪 5 列,再把另外 5 個 bit 設成 1,位圖永遠只有一張,不會變成兩張。
(100 萬個 bit 聽起來很大,實際上用 Roaring bitmap 壓縮,只刪了幾千列的話位圖也才幾 KB。存在一個叫 Puffin 的檔案格式裡,Puffin 是 Iceberg 專門給這種 sidecar 用的容器,deletion vector 是它的第一個 use case。)
於是同一句查詢:
v_delete 直接就是那張位圖,不用另外算代價是每次刪除都要把整張位圖讀回來、改一個 bit、整個 Puffin 檔重寫,寫端變貴了。
Delete file 的兩種寫法(存位置、存主鍵),成本分別壓在寫端跟讀端,可是不管哪一種,delete artifact 都會越積越多,讀的成本會跟著檔案數量一起長。Deletion vector 用「一個資料檔配一張位圖」把數量鎖死成 1,讀的成本變成常數,代價是寫的時候要重寫整張位圖。
因為讀比寫多幾百倍(分析型場景幾乎都這樣),這筆交易划算,所以 Iceberg V3 換成了這條路。Delta Lake 2020 更早就這樣做了,Iceberg V3 才跟上。
三種形式的成本對照大概是這樣:
| 形式 | 寫成本 | 讀成本 | delete artifacts / data file | 適合場景 | 規格 |
|---|---|---|---|---|---|
| Positional delete | 中(要查位置) | 中(sorted merge) | N 份,隨改動累積 | Batch DELETE、小改動 | V2 |
| Equality delete | 低(只要主鍵) | 高(anti-join) | N 份,隨改動累積 | CDC upsert、streaming | V2 |
| Deletion vector | 中偏高(重寫位圖) | 低(bitwise AND) | 1 份 Puffin | 讀多寫少、V3 導入後預設 | V3 |
儲存格式這一層到今天收尾
D9 剖 Parquet 一個檔內部結構、D10 拉出 Iceberg 四層漏斗、D11 講晚物化怎麼在 page 內剪列、D12 講不可變檔案怎麼表達刪除,這四天走完的是「拿到一句 SQL 之後,怎麼從幾百萬個檔剪到只吐出該吐的那幾列」。
回頭看整條漏斗,D10 靠 metadata 樹剪 data file、D9 靠 footer min/max 剪 row group、D11 靠 selection vector 剪 page 內的列、D12 再減掉 delete 標記的列,同一個「證明一定沒有就跳過」的思想被套了四次,每次尺度縮一級,同構重複。
明天 D13 進 DataFusion 引擎這一層,會問一個聽起來很無聊但決定了整個系列命題的問題,緊密整合的查詢引擎(Spark、Presto、DuckDB)已經一堆了,為什麼還要再做一個「可嵌入的模組化引擎」?如果 native 執行是必勝的,那 Spark 這種 JVM 引擎應該早就被淘汰,實際上 EDBT 那篇論文量到某些查詢 vanilla Spark 還贏 native 引擎,這件事就已經很問號了。
留一個沒有標準答案的問題,streaming ingest(每分鐘寫一次 delete)+ heavy analytical read(每秒幾百次讀)的表,要選 equality delete 還是 deletion vector?Equality delete 寫最便宜,來源給你主鍵就寫,但讀時每一列都要 anti-join;deletion vector 讀最便宜,bitwise AND 就完事,但每次寫都要把整張位圖讀回來 merge 再重寫,streaming 頻率越高寫放大越大。答案不只一種,看你 workload 兩邊的比例、看你的 compaction 排程能不能追上、看你能容忍多少讀 latency 換到多少寫 throughput。
那就明天見~